🌎 rns.kin.earth git 🌍
RNS/Interfaces/util/TransmitBuffer.py 9878b1837726109a3a9a5b4e59db067a3e40ca7a (9878b183) Text, 6.77 KB
T8b949e# Reticulum License
T8b949e#
T8b949e# Copyright (c) 2026 Mark Qvist
T8b949e#
T8b949e# Permission is hereby granted, free of charge, to any person obtaining a copy
T8b949e# of this software and associated documentation files (the "Software"), to deal
T8b949e# in the Software without restriction, including without limitation the rights
T8b949e# to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
T8b949e# copies of the Software, and to permit persons to whom the Software is
T8b949e# furnished to do so, subject to the following conditions:
T8b949e#
T8b949e# - The Software shall not be used in any kind of system which includes amongst
T8b949e# its functions the ability to purposefully do harm to human beings.
T8b949e#
T8b949e# - The Software shall not be used, directly or indirectly, in the creation of
T8b949e# an artificial intelligence, machine learning or language model training
T8b949e# dataset, including but not limited to any use that contributes to the
T8b949e# training or development of such a model or algorithm.
T8b949e#
T8b949e# - The above copyright notice and this permission notice shall be included in
T8b949e# all copies or substantial portions of the Software.
T8b949e#
T8b949e# THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
T8b949e# IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
T8b949e# FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
T8b949e# AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
T8b949e# LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
T8b949e# OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
T8b949e# SOFTWARE.
Tff7b72from T7ee787collections Tff7b72import Te6edf3deque
Tff7b72import T7ee787threading
T8b949e# Zero-copy coalesced transmit buffer for arbitrary
T8b949e# frame sizes. It's even almost 100% lock-free.
Tff7b72class T56d364TransmitBufferTb4b4b4:
Te6edf3COALESCE_TARGET Tff7b72= T79c0ff65536
Tff7b72def Tff7b72__init__Tb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4:
Tff7b72selfTff7b72.Td2a8ff_tail_lock Tff7b72= Te6edf3threadingTff7b72.Td2a8ffLockTb4b4b4(Tb4b4b4)
T8b949e# Producer-owned state
Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72= Tff7b72None
Tff7b72selfTff7b72.Td2a8ff_cur_frames Tff7b72= T79c0ff0
T8b949e# Consumer-owned state
Tff7b72selfTff7b72.Td2a8ff_head_offset Tff7b72= T79c0ff0
T8b949e# Shared chunk queue. Producer appends right, consumer
T8b949e# pops from left. Appended chunks are immutable.
Tff7b72selfTff7b72.Td2a8ff_chunks Tff7b72= Te6edf3dequeTb4b4b4(Tb4b4b4)
T8b949e# Accounting counters
T8b949e# Producer-owned
Tff7b72selfTff7b72.Td2a8ff_tx_total Tff7b72= T79c0ff0 T8b949e# Total bytes appended
Tff7b72selfTff7b72.Td2a8ff_tx_visible Tff7b72= T79c0ff0 T8b949e# Total bytes pushed to the chunk queue
Tff7b72selfTff7b72.Td2a8ff_tx_frames Tff7b72= T79c0ff0 T8b949e# Total frames appended
T8b949e# Consumer-owned
Tff7b72selfTff7b72.Td2a8ff_tx_sent Tff7b72= T79c0ff0 T8b949e# Total bytes sent to the kernel
Tff7b72selfTff7b72.Td2a8ff_tx_frames_sent Tff7b72= T79c0ff0 T8b949e# Total frames fully sent
T8b949e# Queue a finalized and complete on-wire frame.
T8b949e# Executed by producer side only.
Tff7b72def Td2a8ffappendTb4b4b4(Tff7b72selfTb4b4b4, Te6edf3frameTb4b4b4, Te6edf3limitTff7b72=Tff7b72NoneTb4b4b4)Tb4b4b4:
Tff7b72with Tff7b72selfTff7b72.Td2a8ff_tail_lockTb4b4b4:
Tff7b72if Te6edf3limit Tff7b72is Tff7b72not Tff7b72None Tff7b72and Tb4b4b4(Tff7b72selfTff7b72.Td2a8ff_tx_total Tff7b72- Tff7b72selfTff7b72.Td2a8ff_tx_sentTb4b4b4) Tff7b72+ Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4) Tff7b72> Te6edf3limitTb4b4b4: Tff7b72return Tff7b72False
T8b949e# Allocate separate chunks to frames
T8b949e# larger than the coalescing target.
Tff7b72if Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4) Tff7b72>Tff7b72= Tff7b72selfTff7b72.Td2a8ffCOALESCE_TARGETTb4b4b4:
Tff7b72selfTff7b72.Td2a8ff_flushTb4b4b4(Tb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_chunksTff7b72.Td2a8ffappendTb4b4b4(Tb4b4b4(Te6edf3frameTb4b4b4, T79c0ff1Tb4b4b4)Tb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_visible Tff7b72+Tff7b72= Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_total Tff7b72+Tff7b72= Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_frames Tff7b72+Tff7b72= T79c0ff1
Tff7b72elseTb4b4b4:
Tff7b72if Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72is Tff7b72not Tff7b72None Tff7b72and Tffa657lenTb4b4b4(Tff7b72selfTff7b72.Td2a8ff_curTb4b4b4) Tff7b72+ Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4) Tff7b72> Tff7b72selfTff7b72.Td2a8ffCOALESCE_TARGETTb4b4b4: Tff7b72selfTff7b72.Td2a8ff_flushTb4b4b4(Tb4b4b4)
Tff7b72if Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72is Tff7b72NoneTb4b4b4:
Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72= Tffa657bytearrayTb4b4b4(Tb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_cur_frames Tff7b72= T79c0ff0
Tff7b72selfTff7b72.Td2a8ff_curTff7b72.Td2a8ffextendTb4b4b4(Te6edf3frameTb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_cur_frames Tff7b72+Tff7b72= T79c0ff1
Tff7b72selfTff7b72.Td2a8ff_tx_total Tff7b72+Tff7b72= Tffa657lenTb4b4b4(Te6edf3frameTb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_frames Tff7b72+Tff7b72= T79c0ff1
Tff7b72if Tff7b72not Tff7b72selfTff7b72.Td2a8ff_chunksTb4b4b4:
T8b949e# Queue is idle. Make the frame visible immediately so
T8b949e# sparse traffic does not incur coalescing latency.
Tff7b72selfTff7b72.Td2a8ff_flushTb4b4b4(Tb4b4b4)
Tff7b72return Tff7b72True
Tff7b72def Td2a8ff_flushTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4:
Tff7b72if Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72is Tff7b72not Tff7b72None Tff7b72and Tffa657lenTb4b4b4(Tff7b72selfTff7b72.Td2a8ff_curTb4b4b4) Tff7b72> T79c0ff0Tb4b4b4:
Tff7b72selfTff7b72.Td2a8ff_chunksTff7b72.Td2a8ffappendTb4b4b4(Tb4b4b4(Tff7b72selfTff7b72.Td2a8ff_curTb4b4b4, Tff7b72selfTff7b72.Td2a8ff_cur_framesTb4b4b4)Tb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_visible Tff7b72+Tff7b72= Tffa657lenTb4b4b4(Tff7b72selfTff7b72.Td2a8ff_curTb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_cur Tff7b72= Tff7b72None
Tff7b72selfTff7b72.Td2a8ff_cur_frames Tff7b72= T79c0ff0
T8b949e# Make producer-side coalescing tail visible to the consumer.
T8b949e# Called by producer when it has finished producing a burst,
T8b949e# and wants to make the remaining bytes transmittable. Safe to
T8b949e# call after every append.
Tff7b72def Td2a8ffflushTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4:
Tff7b72with Tff7b72selfTff7b72.Td2a8ff_tail_lockTb4b4b4: Tff7b72selfTff7b72.Td2a8ff_flushTb4b4b4(Tb4b4b4)
T8b949e# Drain buffered frames to a socket. Called by consumer. If the
T8b949e# chunk queue is drained completely, any potential producer-side
T8b949e# coalescing tail is released under the tail lock, and drained in
T8b949e# the same cycle. If you possess wizardry for doing this without
T8b949e# the lock, I am all ears.
Tff7b72def Td2a8ffdrain_toTb4b4b4(Tff7b72selfTb4b4b4, Te6edf3socketTb4b4b4)Tb4b4b4:
T8b949e# TODO: At *some* point (no, not now), replace with a
T8b949e# mechanism for send_msg to immediately drain all
T8b949e# chunks directly to the kernel.
Te6edf3total Tff7b72= T79c0ff0
Te6edf3whole Tff7b72= Tff7b72True
Tff7b72while Te6edf3wholeTb4b4b4:
Tff7b72if Tff7b72not Tff7b72selfTff7b72.Td2a8ff_chunksTb4b4b4:
Tff7b72with Tff7b72selfTff7b72.Td2a8ff_tail_lockTb4b4b4: Tff7b72selfTff7b72.Td2a8ff_flushTb4b4b4(Tb4b4b4)
Tff7b72if Tff7b72not Tff7b72selfTff7b72.Td2a8ff_chunksTb4b4b4: Tff7b72break
Te6edf3chunkTb4b4b4, Te6edf3_ Tff7b72= Tff7b72selfTff7b72.Td2a8ff_chunksTb4b4b4[T79c0ff0Tb4b4b4]
Te6edf3remaining Tff7b72= Tffa657lenTb4b4b4(Te6edf3chunkTb4b4b4) Tff7b72- Tff7b72selfTff7b72.Td2a8ff_head_offset
Tff7b72if Te6edf3remaining Tff7b72<Tff7b72= T79c0ff0Tb4b4b4:
Tff7b72selfTff7b72.Td2a8ff_pop_headTb4b4b4(Tb4b4b4)
Tff7b72continue
Tff7b72tryTb4b4b4: Te6edf3written Tff7b72= Te6edf3socketTff7b72.Td2a8ffsendTb4b4b4(Tffa657memoryviewTb4b4b4(Te6edf3chunkTb4b4b4)Tb4b4b4[Tff7b72selfTff7b72.Td2a8ff_head_offsetTb4b4b4:Tb4b4b4]Tb4b4b4)
Tff7b72except Tf85149BlockingIOErrorTb4b4b4: Tff7b72break
Tff7b72except Tf85149InterruptedErrorTb4b4b4: Tff7b72continue
Tff7b72if Te6edf3written Tff7b72<Tff7b72= T79c0ff0Tb4b4b4: Tff7b72break
Te6edf3total Tff7b72+Tff7b72= Te6edf3written
Tff7b72selfTff7b72.Td2a8ff_head_offset Tff7b72+Tff7b72= Te6edf3written
Tff7b72selfTff7b72.Td2a8ff_tx_sent Tff7b72+Tff7b72= Te6edf3written
Tff7b72if Tff7b72selfTff7b72.Td2a8ff_head_offset Tff7b72>Tff7b72= Tffa657lenTb4b4b4(Te6edf3chunkTb4b4b4)Tb4b4b4: Tff7b72selfTff7b72.Td2a8ff_pop_headTb4b4b4(Tb4b4b4)
Tff7b72elseTb4b4b4: Te6edf3whole Tff7b72= Tff7b72False
Tff7b72return Te6edf3total
Tff7b72def Td2a8ff_pop_headTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4:
Te6edf3chunkTb4b4b4, Te6edf3frames Tff7b72= Tff7b72selfTff7b72.Td2a8ff_chunksTff7b72.Td2a8ffpopleftTb4b4b4(Tb4b4b4)
Tff7b72selfTff7b72.Td2a8ff_tx_frames_sent Tff7b72+Tff7b72= Te6edf3frames
Tff7b72selfTff7b72.Td2a8ff_head_offset Tff7b72= T79c0ff0
T8b949e# Total number of bytes currently buffered, including the
T8b949e# producer's in-progress coalescing chunk.
Tff7b72def Tff7b72__len__Tb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4: Tff7b72return Tff7b72selfTff7b72.Td2a8ff_tx_totalTff7b72-Tff7b72selfTff7b72.Td2a8ff_tx_sent
T8b949e# The number of bytes actually ready to send, and will be
T8b949e# visible to the consumer, that is, what can be drained.
Tf0883e@property
Tff7b72def Td2a8ffsendableTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4: Tff7b72return Tff7b72selfTff7b72.Td2a8ff_tx_visibleTff7b72-Tff7b72selfTff7b72.Td2a8ff_tx_sent
T8b949e# The number of complete frames currently buffere, including
T8b949e# ones contained in a potential in-progress coalescing chunk.
Tf0883e@property
Tff7b72def Td2a8ffframes_bufferedTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4: Tff7b72return Tff7b72selfTff7b72.Td2a8ff_tx_framesTff7b72-Tff7b72selfTff7b72.Td2a8ff_tx_frames_sent
T8b949e# Number of chunks currently queued and visible to the consumer
Tf0883e@property
Tff7b72def Td2a8ffchunks_bufferedTb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4: Tff7b72return Tffa657lenTb4b4b4(Tff7b72selfTff7b72.Td2a8ff_chunksTb4b4b4)
Tff7b72def Tff7b72__repr__Tb4b4b4(Tff7b72selfTb4b4b4)Tb4b4b4: Tff7b72return Ta5d6fffTa5d6ff"Ta5d6ff<TransmitBuffer buffered=Tffd700{Tffa657lenTb4b4b4(Tff7b72selfTb4b4b4)Tffd700}Ta5d6ff sendable=Tffd700{Tff7b72selfTff7b72.Td2a8ffsendableTffd700}Ta5d6ff frames=Tffd700{Tff7b72selfTff7b72.Td2a8ffframes_bufferedTffd700}Ta5d6ff chunks=Tffd700{Tff7b72selfTff7b72.Td2a8ffchunks_bufferedTffd700}Ta5d6ff>Ta5d6ff"
Served by rngit 1.5.2 - Generated in 0.08s